feat(token-capture): add Megatron worker capture backend - #2823
Conversation
f344be6 to
77950f8
Compare
Add the Megatron worker capture adapter and consolidate backend lifecycle handling. Preserve deferred-call custody and multi-turn lineage, keep capture metadata nested, strip compact prompt token echoes, and cover the shared backend lifecycle with parameterized tests. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
77950f8 to
5f4ddab
Compare
pthombre
left a comment
There was a problem hiding this comment.
1. Add a Gym-owned Megatron extraction adapter
Type: Coordinated Gym/RL design request
Could we add a dependency-light MegatronCaptureAdapter in Gym (nemo_gym/token_id_capture/adapters/megatron.py, mirroring adapters/vllm.py) and have RL #4129 use complete_call_from_response()?
Currently, TQMegatronTokenStager._stage_admitted() performs field extraction and numeric conversion outside Gym's extraction-failure boundary. Missing fields or conversion errors reach the stager's broad exception handler, which logs and returns None. Gym consequently records worker_response_missing_commit_coordinates, whereas equivalent vLLM extraction failures produce explicit capture_failed coordinates and worker_capture_failed.
Moving extraction under the shared wrapper would preserve that distinction and give Megatron payload extraction Gym-owned tests. It also gives Megatron an extract_extras seat; today the stager never passes extras, so routed-experts material is silently absent on MInf. TQ access and policy-epoch handling can remain in RL. Admission/lifecycle errors raised by begin_call still need explicit handling in the consumer; adding the adapter alone will not address those.
One protocol change is needed: CaptureAdapter in staging/protocols.py types response_payload as dict[str, Any], while MInf hands an offloaded payload object. Widening to Any is a runtime no-op since the Protocol is not runtime_checkable.
I would like the adapter file in this PR; the RL-side switch to complete_call_from_response() can follow in #4129.
Resolve the vLLM model conflict with NVIDIA-NeMo#3382 (buffered SSE with external token staging) by porting its stream guards into the capture handler path and adapting the new streaming test to wrap the handler's finalize_response. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
Co-authored-by: Ananth Subramaniam <ananth.subramaniam@gmail.com> Signed-off-by: Laura Dang <lauradang.2000@gmail.com>
| assert "ng_capture" not in outbound | ||
| assert "return_tokens_as_token_ids" not in outbound | ||
|
|
||
| async def test_megatron_capture_handler_finalizes_through_chat_completions(self, monkeypatch: MonkeyPatch) -> None: |
There was a problem hiding this comment.
(for later) @pthombre we should have an end to end test with megatron and vllm in gym that can test token capture since current tests cannot detect mismatches in:
offload_params.ng_capturepropagation through the actual Megatron Inference endpoint- Worker staging before HTTP acknowledgement
ng_commit_coordsplacement- Multi-turn prefix-chain resolution
- Response conversion and stripping across Chat, Responses, Messages, and synthesized SSE
There was a problem hiding this comment.
Agreed. Proposing a follow-up issue scoped to exactly these five mismatch classes (offload_params.ng_capture propagation, worker staging before HTTP ack, ng_commit_coords placement, multi-turn prefix-chain resolution, and response conversion/stripping across Chat, Responses, Messages, and synthesized SSE) with a real Megatron Inference and vLLM endpoint. Leaving this thread open as the tracker until that issue is filed.
|
|
||
| def extract_prompt_ids(self, response_payload: dict[str, Any]) -> list[int]: | ||
| """Return the exact prompt token IDs used for generation.""" | ||
| def extract_prompt_ids(self, response_payload: Any) -> list[int]: |
There was a problem hiding this comment.
@pthombre potentially for later, we could better type this with a generic CaptureAdapter[PayloadT] to get some static typechecking help here in addition to relying on the adapter implementation to refine the type
There was a problem hiding this comment.
Agreed this is worth doing as a follow-up. Note the scope is a bit wider than the protocol: install_capture takes adapter: CaptureAdapter | None, and the worker consumes the payload generically, so a PayloadT parameter only gives static help once the install path and worker are parametrized too. Leaving this thread open as the tracker until a follow-up issue exists.
An unpoisoned receipt must name both terminal_model_call_id and terminal_selection and carry no failure_reason. An unset terminal_selection is reserved for receipts where no attribution stage ran: such a receipt names no terminal and must be poisoned with a reason. A poisoned receipt may still retain its terminal and method, because attribution can succeed before a later call poisons the rollout. The linearize compatibility wrapper now reports missing_terminal itself instead of constructing a receipt the schema rejects. Signed-off-by: Laura Dang <laurad@nvidia.com> Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
…apture input Validate element types in the Megatron adapter before conversion: token id fields must contain only plain integers and log-probability fields only numbers, so a float, string, or bool element poisons the call instead of being silently coerced into a plausible id. Reject megatron_worker capture requests whose messages carry non-text content parts. The adapter stages no media geometry yet, so a multimodal prompt would commit rows misaligned with the expanded engine prompt. Spell out Megatron Inference instead of the MInf shorthand. Signed-off-by: Laura Dang <laurad@nvidia.com> Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
…fig check Signed-off-by: Laura Dang <laurad@nvidia.com> Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test b601670 |
|
🔄 Merge queue validation started! You can track the progress here: https://github.com/NVIDIA-NeMo/Gym/actions/runs/35408060167 |
Move 3rdparty/Gym-workspace/Gym from the laurad/megatron-adapter fork branch (37dc751f4) to NVIDIA-NeMo/Gym main at 9fc05c0ff, the squash merge of NVIDIA-NeMo/Gym#2823 (Megatron worker capture backend). The pin now lives on the submodule's configured upstream and branch, so CI and fresh clones can fetch it. Every nemo_gym symbol imported from nemo_rl/ and tests/ resolves at the new pin, and Gym's pyproject is identical between the two pins, so uv.lock is unaffected. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
Summary
megatron_workerexternal staging backend so Megatron-Inference (MInf) serving workers can participate in framework-owned token captureexternal_capture.py) with a shared prepare/finalize lifecycleng_commit_coords, and Gym records a token-freeCallRecordonly after that acknowledgementreturn_tokenized_datainstead ofreturn_tokens_as_token_ids); both strip token IDs, logprobs, and coordinates off the agent hopbegin_callaccept an explicitweight_versionso the MInf worker can stamp the epoch from engine-finished metadataTesting
pre-commit run --fileson all changed working-tree filesBaseline note
The complete
responses_api_models/vllm_model/tests/test_app.pyrun has 15 response-ID expectation failures (chtcmpl-123versusresp_123). A representative failure reproduces unchanged onpthombre/tq-tokidcap-capture, so it is not introduced by this branch. The remaining red CI checks (lint, copyright,test_server_utils) come from files this PR does not touch, and the shard 2/6 failures are an apt mirror error.Dependencies
pthombre/tq-tokidcap-capture, notmain.🤖 Generated with Claude Code